Dedup duplicate join notifies + reject joins on dead broadcaster - #21
Merged
Merged
Conversation
Live deploy at 1080p60 hit two compounding bugs the moment a viewer clicked Join: 1. The TS6 server retransmitted ``notifyjoinstreamrequest`` 9 times inside 16 ms (server quirk; we ACK at the TS3 layer but it resent anyway). The publisher's existing dedup looked at ``self._viewers.pop(...)`` *before* the slow path - but the slow path takes tens of milliseconds (createOffer + ICE gathering) and the viewer slot only lands in ``_viewers`` after that work completes. So all 9 retransmits raced past the dedup, each spawning a fresh peer connection + broadcaster subscription before any of them could populate the dict. Nine concurrent ICE gatherings + nine encoder subscriptions blew RSS to 1.3 GB and OOM-killed ffmpeg/x11grab. 2. Once the broadcaster's source raised MediaStreamError on the OOM-killed capture, the pump task exited but ``subscribe()`` kept handing out queues that nothing would ever fill. Newly joining viewers got an offer that pointed at a dead media pipeline; their UI hung on "connecting" forever. Fixes: * ``StreamPublisher._joining: set[int]`` lock-protected, populated *before* any expensive work (offer creation, ICE, broadcaster subscribe) and discarded in a finally block. Duplicate notifies for the same clid log ``stream_publisher.duplicate_join_dropped`` and bail out without touching aiortc or the broadcaster. * ``VideoBroadcaster.is_alive`` (False after the pump task exits for any reason). When false, ``_handle_viewer_join`` sends ``respondjoinstreamrequest decision=0`` so the TS6 client UI fails the join cleanly rather than hanging. The test mock for ``RTCPeerConnection.createOffer`` / ``setLocalDescription`` now actually yields via ``await asyncio.sleep(0)``. Without that, the in-flight dedup looked effective only because the mocked path ran to completion synchronously - in real aiortc the ICE-gathering yield would let parallel tasks observe each other's ``_joining`` state, which is exactly what we want to assert.
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Symptom
Live deploy at 1080p60 / 6 Mbit / viewer_limit=5: bot connects fine, stream allocates, then on the very first viewer click the bot goes into a tailspin. The log told the story:
The TS6 server retransmitted the same
notifyjoinstreamrequest9 times inside 16 ms. Our existing dedup checkedself._viewers.pop(...)before the slow path (createOffer+ ICE gathering = tens of ms), but the viewer slot only lands in_viewersafter that slow path completes. All 9 retransmits raced past the dedup, each spawning a fresh peer connection + broadcaster subscription. The 9× concurrent ICE gatherings + encoder subscriptions blew RSS to 1.3 GB and OOM-killed the capture pipeline. After that, the broadcaster's pump task exited butsubscribe()kept handing out queues nothing would ever fill — so subsequent viewers' UI hung on "connecting" forever (= "kommt nicht wieder").Fixes
StreamPublisher._joining— a lock-protectedset[int]populated before any expensive work (offer creation, ICE, broadcaster subscribe) and discarded in a finally block. Duplicate notifies for the same clid logstream_publisher.duplicate_join_droppedand bail out cheaply without touching aiortc or the broadcaster.VideoBroadcaster.is_alive—Falseafter the pump task exits for any reason (sourceMediaStreamError, encoder crash). When false,_handle_viewer_joinsendsrespondjoinstreamrequest decision=0and exits — the TS6 client UI fails the join cleanly instead of hanging on "connecting".Test plan
pytest— 250 passed, 1 skipped.test_duplicate_join_for_same_clid_only_creates_one_pcbursts 9 duplicate notifies and asserts exactly one PC, one broadcaster subscription, one audio relay subscription.test_join_refused_when_broadcaster_source_deadflipsis_alive=Falseand asserts decision=0 in the response with no PC, no subscriptions.test_is_alive_*covers the broadcaster's lifecycle flag.createOffer/setLocalDescriptionnowawait asyncio.sleep(0)s — without it the in-flight dedup looked effective only because the mocked path ran to completion synchronously, but real aiortc's ICE-gathering yield is what lets parallel tasks observe each other's_joiningstate.mypy --strictclean.ruff checkclean.stream_publisher.duplicate_join_droppedlogs for the 8 retransmits, no RSS spike, no OOM.Operator notes
Even with this fix, 1080p60 is still asking a lot from libvpx-realtime on a small VPS. If the bot stays alive after this PR but the encoder still can't keep up (frames pile up, RSS climbs), drop to 1080p30:
Or 720p60:
https://claude.ai/code/session_016DuCjRJK995Tj9aDhhB9at
Generated by Claude Code